--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
Commit 4fb50f412ed86553963c0666032588b402580278
Parents : 4418f91
Author : Sudo-Ivan <ivan@quad4.io>
Signature : Signature validation error
Date : 2026-03-06T03:54:36-06:00
Improve message retrieval and database operations
- Implemented pagination for message export in the maintenance API to handle large datasets efficiently.
- Set default limit values for message queries in MessageHandler and AnnounceDAO to improve consistency.
- Added a busy timeout pragma in the database initialization to manage lock contention.
- Refactored batch insert operations in MessageDAO for improved performance and readability.
Changes
7 files changed, 84 insertions(+), 42 deletions(-)
Diff
diff --git a/meshchatx/meshchat.py b/meshchatx/meshchat.py
index 9a410b7e..bdff6801 100644
--- a/meshchatx/meshchat.py
+++ b/meshchatx/meshchat.py
@@ -4594,9 +4594,17 @@ class ReticulumMeshChat:
# maintenance - export messages
@routes.get("/api/v1/maintenance/messages/export")
async def maintenance_export_messages(request):
- messages = self.database.messages.get_all_lxmf_messages()
- # Convert sqlite3.Row to dict if necessary
- messages_list = [dict(m) for m in messages]
+ messages_list = []
+ page_size = 5000
+ offset = 0
+ while True:
+ page = self.database.messages.get_all_lxmf_messages(
+ limit=page_size, offset=offset
+ )
+ messages_list.extend(dict(m) for m in page)
+ if len(page) < page_size:
+ break
+ offset += page_size
return web.json_response({"messages": messages_list})
# maintenance - import messages
diff --git a/meshchatx/src/backend/announce_manager.py b/meshchatx/src/backend/announce_manager.py
index 6c2e487f..1d04f62d 100644
--- a/meshchatx/src/backend/announce_manager.py
+++ b/meshchatx/src/backend/announce_manager.py
@@ -49,7 +49,7 @@ class AnnounceManager:
destination_hash=None,
query=None,
blocked_identity_hashes=None,
- limit=None,
+ limit=500,
offset=0,
):
sql = """
diff --git a/meshchatx/src/backend/database/__init__.py b/meshchatx/src/backend/database/__init__.py
index c5c3db96..5d37e947 100644
--- a/meshchatx/src/backend/database/__init__.py
+++ b/meshchatx/src/backend/database/__init__.py
@@ -67,6 +67,7 @@ class Database:
self.execute_sql("PRAGMA temp_store=MEMORY")
self.execute_sql("PRAGMA cache_size=-8000") # 8 MB
self.execute_sql("PRAGMA mmap_size=67108864") # 64 MB
+ self.execute_sql("PRAGMA busy_timeout=5000") # 5 s wait on lock contention
except Exception as exc:
print(f"SQLite pragma setup failed: {exc}")
diff --git a/meshchatx/src/backend/database/announces.py b/meshchatx/src/backend/database/announces.py
index a1b626e8..c9172ef1 100644
--- a/meshchatx/src/backend/database/announces.py
+++ b/meshchatx/src/backend/database/announces.py
@@ -69,7 +69,7 @@ class AnnounceDAO:
search_term=None,
identity_hash=None,
destination_hash=None,
- limit=None,
+ limit=500,
offset=0,
):
query = "SELECT * FROM announces WHERE 1=1"
diff --git a/meshchatx/src/backend/database/messages.py b/meshchatx/src/backend/database/messages.py
index 4548397f..e7e3910b 100644
--- a/meshchatx/src/backend/database/messages.py
+++ b/meshchatx/src/backend/database/messages.py
@@ -80,11 +80,15 @@ class MessageDAO:
)
def delete_all_lxmf_messages(self):
- self.provider.execute("DELETE FROM lxmf_messages")
- self.provider.execute("DELETE FROM lxmf_conversation_read_state")
+ with self.provider:
+ self.provider.execute("DELETE FROM lxmf_messages")
+ self.provider.execute("DELETE FROM lxmf_conversation_read_state")
- def get_all_lxmf_messages(self):
- return self.provider.fetchall("SELECT * FROM lxmf_messages")
+ def get_all_lxmf_messages(self, limit=5000, offset=0):
+ return self.provider.fetchall(
+ "SELECT * FROM lxmf_messages ORDER BY id LIMIT ? OFFSET ?",
+ (limit, offset),
+ )
def count_lxmf_messages(self):
row = self.provider.fetchone("SELECT COUNT(*) AS count FROM lxmf_messages")
@@ -103,10 +107,14 @@ class MessageDAO:
(destination_hash, limit, offset),
)
+ CONVERSATION_LIST_COLUMNS = (
+ "m1.id, m1.hash, m1.source_hash, m1.destination_hash, m1.peer_hash, "
+ "m1.state, m1.is_incoming, m1.title, m1.timestamp, m1.created_at, m1.updated_at"
+ )
+
def get_conversations(self):
- # Optimized using peer_hash column
- query = """
- SELECT m1.* FROM lxmf_messages m1
+ query = f"""
+ SELECT {self.CONVERSATION_LIST_COLUMNS} FROM lxmf_messages m1
INNER JOIN (
SELECT peer_hash, MAX(timestamp) as max_ts
FROM lxmf_messages
@@ -115,7 +123,7 @@ class MessageDAO:
) m2 ON m1.peer_hash = m2.peer_hash AND m1.timestamp = m2.max_ts
GROUP BY m1.peer_hash
ORDER BY m1.timestamp DESC
- """
+ """ # noqa: S608
return self.provider.fetchall(query)
def mark_conversation_as_read(self, destination_hash):
@@ -136,17 +144,16 @@ class MessageDAO:
return
now = datetime.now(UTC).isoformat()
with self.provider:
- for destination_hash in destination_hashes:
- self.provider.execute(
- """
- INSERT INTO lxmf_conversation_read_state (destination_hash, last_read_at, created_at, updated_at)
- VALUES (?, ?, ?, ?)
- ON CONFLICT(destination_hash) DO UPDATE SET
- last_read_at = EXCLUDED.last_read_at,
- updated_at = EXCLUDED.updated_at
- """,
- (destination_hash, now, now, now),
- )
+ self.provider.executemany(
+ """
+ INSERT INTO lxmf_conversation_read_state (destination_hash, last_read_at, created_at, updated_at)
+ VALUES (?, ?, ?, ?)
+ ON CONFLICT(destination_hash) DO UPDATE SET
+ last_read_at = EXCLUDED.last_read_at,
+ updated_at = EXCLUDED.updated_at
+ """,
+ [(h, now, now, now) for h in destination_hashes],
+ )
def is_conversation_unread(self, destination_hash):
row = self.provider.fetchone(
@@ -312,17 +319,16 @@ class MessageDAO:
now = datetime.now(UTC).isoformat()
if destination_hashes:
with self.provider:
- for destination_hash in destination_hashes:
- self.provider.execute(
- """
- INSERT INTO notification_viewed_state (destination_hash, last_viewed_at, created_at, updated_at)
- VALUES (?, ?, ?, ?)
- ON CONFLICT(destination_hash) DO UPDATE SET
- last_viewed_at = EXCLUDED.last_viewed_at,
- updated_at = EXCLUDED.updated_at
- """,
- (destination_hash, now, now, now),
- )
+ self.provider.executemany(
+ """
+ INSERT INTO notification_viewed_state (destination_hash, last_viewed_at, created_at, updated_at)
+ VALUES (?, ?, ?, ?)
+ ON CONFLICT(destination_hash) DO UPDATE SET
+ last_viewed_at = EXCLUDED.last_viewed_at,
+ updated_at = EXCLUDED.updated_at
+ """,
+ [(h, now, now, now) for h in destination_hashes],
+ )
else:
# mark all conversations as viewed
self.provider.execute(
@@ -399,9 +405,27 @@ class MessageDAO:
)
def move_conversations_to_folder(self, peer_hashes, folder_id):
+ if not peer_hashes:
+ return
+ now = datetime.now(UTC).isoformat()
with self.provider:
- for peer_hash in peer_hashes:
- self.move_conversation_to_folder(peer_hash, folder_id)
+ if folder_id is None:
+ placeholders = ", ".join(["?"] * len(peer_hashes))
+ self.provider.execute(
+ f"DELETE FROM lxmf_conversation_folders WHERE peer_hash IN ({placeholders})", # noqa: S608
+ tuple(peer_hashes),
+ )
+ else:
+ self.provider.executemany(
+ """
+ INSERT INTO lxmf_conversation_folders (peer_hash, folder_id, created_at, updated_at)
+ VALUES (?, ?, ?, ?)
+ ON CONFLICT(peer_hash) DO UPDATE SET
+ folder_id = EXCLUDED.folder_id,
+ updated_at = EXCLUDED.updated_at
+ """,
+ [(h, folder_id, now, now) for h in peer_hashes],
+ )
def get_all_conversation_folders(self):
return self.provider.fetchall("SELECT * FROM lxmf_conversation_folders")
diff --git a/meshchatx/src/backend/database/provider.py b/meshchatx/src/backend/database/provider.py
index 1ac8f690..d46f4bc7 100644
--- a/meshchatx/src/backend/database/provider.py
+++ b/meshchatx/src/backend/database/provider.py
@@ -126,6 +126,11 @@ class DatabaseProvider:
else:
self.commit()
+ def executemany(self, query, params_seq):
+ cursor = self.connection.cursor()
+ cursor.executemany(query, params_seq)
+ return cursor
+
def fetchone(self, query, params=None):
cursor = self.execute(query, params)
row = cursor.fetchone()
diff --git a/meshchatx/src/backend/message_handler.py b/meshchatx/src/backend/message_handler.py
index cf959852..63045035 100644
--- a/meshchatx/src/backend/message_handler.py
+++ b/meshchatx/src/backend/message_handler.py
@@ -41,15 +41,17 @@ class MessageHandler:
[destination_hash],
)
- def search_messages(self, local_hash, search_term):
+ def search_messages(self, local_hash, search_term, limit=500):
like_term = f"%{search_term}%"
query = """
SELECT peer_hash, MAX(timestamp) as max_ts
FROM lxmf_messages
WHERE title LIKE ? OR content LIKE ? OR peer_hash LIKE ?
GROUP BY peer_hash
+ ORDER BY max_ts DESC
+ LIMIT ?
"""
- params = [like_term, like_term, like_term]
+ params = [like_term, like_term, like_term, limit]
return self.db.provider.fetchall(query, params)
def get_conversations(
@@ -60,13 +62,15 @@ class MessageHandler:
filter_failed=False,
filter_has_attachments=False,
folder_id=None,
- limit=None,
+ limit=500,
offset=0,
):
- # Optimized using peer_hash column and JOINs to avoid N+1 queries
query = """
SELECT
- m1.*,
+ m1.id, m1.hash, m1.source_hash, m1.destination_hash,
+ m1.peer_hash, m1.state, m1.progress, m1.is_incoming,
+ m1.title, m1.timestamp, m1.is_spam, m1.reply_to_hash,
+ m1.created_at, m1.updated_at,
a.app_data as peer_app_data,
c.display_name as custom_display_name,
con.custom_image as contact_image,
──────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────